آموزش کامل Apache Kafka در ASP.NET Core

Kafka asp.net core همه سطوح
parsakarimidev.ir

آموزش کامل Apache Kafka در ASP.NET Core

Kafka asp.net coreفریم‌ورک سطح: همه سطوح 1405/07/16

آموزش کامل Apache Kafka در ASP.NET Core

⚡ خلاصه سریع — اصل کاری‌ها

این‌ها مهم‌ترین چیزهایی هستن که اگه هیچی ندونی، فقط اینا رو بدونی کافیه:

  • نکته ۱: Kafka یک Message Broker توزیع‌شده است که پیام‌ها رو به‌صورت جریانی (Stream) ذخیره و مدیریت می‌کنه.
  • نکته ۲: داده‌ها در Topic ذخیره می‌شن و هر Topic می‌تونه به چند Partition تقسیم بشه.
  • نکته ۳: Producer پیام تولید می‌کنه، Consumer پیام مصرف می‌کنه.
  • نکته ۴: Consumer Group باعث می‌شه چند Consumer بار پردازش رو بین خودشون تقسیم کنن.
  • نکته ۵: Offset شماره هر پیام تو Partition هست و Consumer با اون می‌دونه کجا مونده.
  • نکته ۶: پیام‌ها بعد از مصرف حذف نمی‌شن (تا زمان Retention)، این تفاوت اصلی Kafka با RabbitMQ است.
  • نکته ۷: برای کار با Kafka در .NET از پکیج Confluent.Kafka استفاده می‌شه.
  • نکته ۸: Kafka روی Zookeeper یا حالت جدید KRaft اجرا می‌شه.
  • نکته ۹: Kafka داده رو روی دیسک ذخیره می‌کنه و به‌خاطر Sequential I/O خیلی سریعه.
  • نکته ۱۰: Kafka حداقل یه‌بار تحویل (At-least-once) تضمین می‌کنه، نه دقیقاً یه‌بار.

۳ اشتباه رایج که باید ازشون پرهیز کنی:

  1. Offset رو خودت Commit نمی‌کنی و پیام‌ها گم می‌شن یا دوبار پردازش می‌شن.
  2. تعداد Partitionها رو بی‌دلیل زیاد می‌کنی که overhead ایجاد می‌کنه.
  3. خطاهای Deserialization رو تو Consumer هندل نمی‌کنی و کل برنامه هنگ می‌کنه.

اگر فقط ۵ دقیقه وقت داری، این‌ها رو یاد بگیر:

  1. Kafka = Topic + Producer + Consumer (ساده‌ترین تصویر ذهنی)
  2. dotnet add package Confluent.Kafka (پکیج اصلی کار با Kafka در .NET)
  3. یه Producer بساز، پیام بفرست، یه Consumer بساز، پیام رو بخون

۱. مقدمه

این موضوع چیه؟

Apache Kafka یه پلتفرم توزیع‌شده برای Event Streaming است. به زبان ساده‌تر، یه صف پیام (Message Queue) خیلی قوی و سریعه که می‌تونه میلیون‌ها پیام در ثانیه رو پردازش کنه.

تصور کن یه رستوران خیلی شلوغ داری:

  • آشپز = Producer (پیام تولید می‌کنه — مثلاً یه سفارش غذا)
  • گارسون = Kafka (پیام رو می‌گیره و تا زمان مصرف نگه می‌داره)
  • میزبان = Consumer (پیام رو دریافت و پردازش می‌کنه)

تفاوت Kafka با صف‌های معمولی اینه که گارسون (Kafka) غذا رو دور نمی‌ریزه! اون نگهش می‌داره. اگه چندتا میزبان (Consumer) بیان، هرکدوم یه بخش از غذاها رو می‌گیرن و اگه یه میزبان غیبت کنه، غذاها هنوز هستن.

چرا باید یادش بگیری؟

دلیل توضیح
سرعت بالا می‌تونه میلیون‌ها پیام در ثانیه پردازش کنه
پایداری (Durability) پیام‌ها روی دیسک ذخیره می‌شن، حتی اگه سرور crash کنه
مقیاس‌پذیری می‌تونی چندتا Broker اضافه کنی و بار رو تقسیم کنی
Replay می‌تونی پیام‌های قدیمی رو دوباره بخونی
تقاضای بازار شرکت‌های بزرگی مثل LinkedIn، Netflix، Uber ازش استفاده می‌کنن

کجا استفاده می‌شه؟

  • سیستم‌های Microservices — ارتباط بین سرویس‌ها بدون وابستگی مستقیم
  • Log Aggregation — جمع‌آوری لاگ از سرورهای مختلف
  • Event Sourcing — ذخیره تغییرات state به‌صورت Event
  • Real-time Analytics — تحلیل داده‌ها در لحظه
  • Stream Processing — پردازش جریان داده‌ها
  • نبود Single Point of Failure — اگه یه سرور بده، بقیه کار می‌کنن

۲. پیش‌نیازها

چه چیزهایی رو باید از قبل بدونی

مهارت سطح مورد نیاز
C# متوسط (LINQ، async/await، Generics)
ASP.NET Core مقدماتی (Minimal API یا Controller)
Docker مقدماتی (run، compose)
Command Line مقدماتی
مفاهیم Messaging آشنایی کلی

ابزارهای لازم

  • .NET SDK 8.0 یا بالاتر — دانلود از سایت مایکروسافت
  • Docker Desktop — برای اجرای Kafka به‌صورت محلی
  • Visual Studio 2022 یا VS Code یا JetBrains Rider
  • kafka-topics CLI — ( داخل Docker هست، نیاز جداگانه نیست )

۳. نصب و راه‌اندازی

گام ۱: اجرای Kafka با Docker Compose

یه فایل به اسم docker-compose.yml بساز:

# docker-compose.yml
version: '3.8'

services:
  kafka:
    image: bitnami/kafka:latest
    ports:
      - "9092:9092"
    environment:
      # KRaft mode (بدون نیاز به Zookeeper)
      KAFKA_CFG_NODE_ID: 0
      KAFKA_CFG_PROCESS_ROLES: controller,broker
      KAFKA_CFG_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093
      KAFKA_CFG_LISTENERS_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT
      KAFKA_CFG_CONTROLLER_QUORUM_VOTERS: 0@kafka:9093
      KAFKA_CFG_CONTROLLER_LISTENER_NAMES: CONTROLLER
      KAFKA_CFG_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
    volumes:
      - kafka_data:/bitnami/kafka

volumes:
  kafka_data:

حالا اجراش کن:

docker-compose up -d

توضیح: این پیکربندی از حالت KRaft استفاده می‌کنه که دیگه نیازی به Zookeeper نداره. Kafka خودش مدیریت Controller رو انجام می‌ده. خیلی ساده‌تر و سبک‌تره.

گام ۲: بررسی سلامت Kafka

docker exec -it <container_id> kafka-topics.sh --list --bootstrap-server localhost:9092

اگه بدون خطا اجرا شد و لیست خالی برگردوند، یعنی Kafka بالا اومده.

گام ۳: ساخت اولین پروژه ASP.NET Core

# ساخت solution
dotnet new sln -n KafkaTutorial

# ساخت Web API project
dotnet new webapi -n KafkaTutorial.Api -o src/KafkaTutorial.Api

# اضافه کردن پروژه به solution
dotnet sln add src/KafkaTutorial.Api

# نصب پکیج Kafka
cd src/KafkaTutorial.Api
dotnet add package Confluent.Kafka

توضیح: پکیج Confluent.Kafka کتابخانه رسمیه که Confluent (شرکت پشت Kafka) برای .NET منتشر می‌کنه. روی librdkafka (کتابخانه C) ساخته شده و performance بالایی داره.


۴. مفاهیم پایه

بیاید قبل از کد، مفاهیم رو با مثال ساده درک کنیم.

۴.۱ Broker

Broker یه سرور Kafka است. اگه ۳ تا سرور Kafka داشته باشی، ۳ تا Broker داری. Brokerها با هم صحبت می‌کنن و داده‌ها رو بین خودشون توزیع می‌کنن.

┌──────────┐  ┌──────────┐  ┌──────────┐
│ Broker 1 │  │ Broker 2 │  │ Broker 3 │
│  :9092   │←→│  :9093   │←→│  :9094   │
└──────────┘  └──────────┘  └──────────┘

۴.۲ Topic

Topic یه دسته‌بندی پیام‌هاست. مثل یه کانال تلگرام — همه پیام‌های مربوط به یه موضوع اونجا قرار می‌گیرن.

مثال: orders، user-events، payment-notifications

۴.۳ Partition

هر Topic به چندتا Partition تقسیم می‌شه. هر Partition یه log مرتب و تغییرناپذیر (Append-only) از پیام‌هاست.

Topic: orders
├── Partition 0: [msg1, msg2, msg5, msg8]
├── Partition 1: [msg3, msg4, msg6, msg9]
└── Partition 2: [msg7, msg10, msg11]

چرا Partition مهمه؟ چون Partition واحد موازی‌سازی است. اگه ۳ Partition داشته باشی، می‌تونی ۳ Consumer موازی داشته باشی که هرکدوم یه Partition رو می‌خونن.

۴.۴ Offset

هر پیام تو یه Partition یه شماره داره به اسم Offset. از ۰ شروع می‌شه و یکی یکی اضافه می‌شه.

Partition 0: [offset 0] [offset 1] [offset 2] [offset 3]

Consumer با ذخیره Offset فعلی، می‌دونه کجا مونده و بعداً از همون‌جا ادامه می‌ده.

۴.۵ Producer

Producer پیام رو به Topic می‌فرسته. خودش تصمیم می‌گیره (یا با یه Key) پیام به کدوم Partition بره.

۴.۶ Consumer

Consumer پیام‌ها رو از Topic می‌خونه. می‌تونه از یه Offset خاص شروع کنه.

۴.۷ Consumer Group

یه گروه از Consumerها که با هم کار می‌کنن و Partitionها رو بین خودشون تقسیم می‌کنن.

Consumer Group: "order-processors"
├── Consumer A ← reads Partition 0
├── Consumer B ← reads Partition 1
└── Consumer C ← reads Partition 2

قانون مهم: اگه تعداد Consumer بیشتر از Partition باشه، بعضی Consumerها هیچ کاری نمی‌کنن (idle می‌شن).

۴.۸ Retention

Kafka پیام‌ها رو برای یه مدت مشخص نگه می‌داره (پیش‌فرض ۷ روز). بعد از اون، خودکار حذف می‌شن.

۴.۹ تجسم کلی

 Producer                 Kafka Cluster                  Consumer
 ┌──────┐     ┌───────────────────────────┐          ┌──────────┐
 │ App  │────→│  Topic: orders            │ ────────→│ Consumer │
 │      │     │  ├── Partition 0          │          │ Group A  │
 └──────┘     │  ├── Partition 1          │          └──────────┘
              │  └── Partition 2          │          ┌──────────┐
              └───────────────────────────┘ ────────→│ Consumer │
                                                    │ Group B  │
                                                    └──────────┘

۵. مثال‌های کد ساده

۵.۱ ساخت Topic با کد

اول یه فولدر Kafka تو پروژه بساز و بعد یه کلاس KafkaConfig:

// Kafka/KafkaConfig.cs
namespace KafkaTutorial.Api.Kafka;

public static class KafkaConfig
{
    public const string BootstrapServers = "localhost:9092";
    public const string TopicName = "orders";
}

توضیح: این کلاس تنظیمات اصلی Kafka رو نگه می‌داره. BootstrapServers آدرس سرور Kafka است. وقتی توسعه محلی می‌کنی، localhost:9092 کافیه.

۵.۲ Producer ساده — ارسال یه پیام

یه کلاس KafkaProducerService بساز:

// Kafka/KafkaProducerService.cs
using Confluent.Kafka;

namespace KafkaTutorial.Api.Kafka;

public class KafkaProducerService
{
    private readonly IProducer<Null, string> _producer;
    private readonly string _topic = KafkaConfig.TopicName;

    public KafkaProducerService()
    {
        var config = new ProducerConfig
        {
            BootstrapServers = KafkaConfig.BootstrapServers
        };

        _producer = new ProducerBuilder<Null, string>(config).Build();
    }

    public async Task SendMessageAsync(string message)
    {
        var result = await _producer.ProduceAsync(
            _topic,
            new Message<Null, string> { Value = message }
        );

        Console.WriteLine(
            $"✅ پیام ارسال شد — " +
            $"Topic: {result.Topic}, " +
            $"Partition: {result.Partition}, " +
            $"Offset: {result.Offset}"
        );
    }
}

توضیح:

  • IProducer<Null, string> یعنی Key پیام null است و Value از نوع string.
  • ProduceAsync پیام رو به Kafka می‌فرسته و منتظر می‌شه تا تأیید بشه.
  • result اطلاعات کامل پیام رو برمی‌گردونه — کدوم Partition، Offset چقدر شده.

۵.۳ استفاده از Producer تو Minimal API

// Program.cs
using KafkaTutorial.Api.Kafka;

var builder = WebApplication.CreateBuilder(args);

builder.Services.AddSingleton<KafkaProducerService>();

var app = builder.Build();

app.MapPost("/orders", async (KafkaProducerService producer, string order) =>
{
    await producer.SendMessageAsync(order);
    return Results.Ok(new { Message = "سفارش ثبت شد", Order = order });
});

app.Run();

توضیح: الان یه endpoint ساختی که پیام رو به‌صورت POST می‌گیره و به Kafka می‌فرسته. تست:

curl -X POST "http://localhost:5000/orders?order=pizza-margherita"

۵.۴ Consumer ساده — خواندن پیام

یه کلاس KafkaConsumerService بساز که به‌صورت Background Service اجرا بشه:

// Kafka/KafkaConsumerService.cs
using Confluent.Kafka;

namespace KafkaTutorial.Api.Kafka;

public class KafkaConsumerService : BackgroundService
{
    private readonly string _topic = KafkaConfig.TopicName;
    private readonly string _groupId = "order_processor_group";
    private readonly string _bootstrapServers = KafkaConfig.BootstrapServers;

    protected override Task ExecuteAsync(CancellationToken stoppingToken)
    {
        return Task.Run(() => StartConsumerLoop(stoppingToken), stoppingToken);
    }

    private void StartConsumerLoop(CancellationToken cancellationToken)
    {
        var config = new ConsumerConfig
        {
            BootstrapServers = _bootstrapServers,
            GroupId = _groupId,
            AutoOffsetReset = AutoOffsetReset.Earliest
        };

        using var consumer = new ConsumerBuilder<Null, string>(config).Build();
        consumer.Subscribe(_topic);

        try
        {
            while (!cancellationToken.IsCancellationRequested)
            {
                var consumeResult = consumer.Consume(cancellationToken);
                Console.WriteLine(
                    $"📩 پیام دریافت شد — " +
                    $"Partition: {consumeResult.Partition}, " +
                    $"Offset: {consumeResult.Offset}, " +
                    $"Value: {consumeResult.Message.Value}"
                );
            }
        }
        catch (OperationCanceledException)
        {
            Console.WriteLine("Consumer متوقف شد.");
        }
        finally
        {
            consumer.Close();
        }
    }
}

توضیح:

  • BackgroundService یعتی این سرویس تو پس‌زمینه اجرا می‌شه.
  • AutoOffsetReset.Earliest یعنی اولین باری که Consumer اجرا می‌شه، از ابتدای Topic بخونه.
  • Consume blocking است — منتظر می‌مونه تا پیام بیاد.
  • consumer.Close() یه Commit خودکار از Offset فعلی انجام می‌ده.

۵.۵ ثبت Consumer تو Program.cs

// Program.cs (بخش‌های اضافه‌شده)
builder.Services.AddHostedService<KafkaConsumerService>();

۵.۶ تست کار

حالا برنامه رو اجرا کن:

dotnet run

یه ترمینال دیگه باز کن و پیام بفرست:

curl -X POST "http://localhost:5000/orders?order=pizza-margherita"
curl -X POST "http://localhost:5000/orders?order=burger"
curl -X POST "http://localhost:5000/orders?order=salad"

تو Console برنامه باید ببینی:

✅ پیام ارسال شد — Topic: orders, Partition: 0, Offset: 0
✅ پیام ارسال شد — Topic: orders, Partition: 0, Offset: 1
✅ پیام ارسال شد — Topic: orders, Partition: 0, Offset: 2
📩 پیام دریافت شد — Partition: 0, Offset: 0, Value: pizza-margherita
📩 پیام دریافت شد — Partition: 0, Offset: 1, Value: burger
📩 پیام دریافت شد — Partition: 0, Offset: 2, Value: salad

۶. مثال واقعی و کاربردی

پروژه کوچک: سیستم سفارش‌گذاری

سناریو: یه رستوران آنلاین داریم. وقتی مشتری سفارش می‌ده:

  1. OrderService یه پیام به Topic orders می‌فرسته
  2. KitchenService (Consumer) سفارش رو دریافت و پردازش می‌کنه
  3. NotificationService (Consumer دیگه) ایمیل تأیید می‌فرسته

گام ۱: مدل داده

// Models/Order.cs
using System.Text.Json.Serialization;

namespace KafkaTutorial.Api.Models;

public class Order
{
    public Guid Id { get; set; } = Guid.NewGuid();
    public string CustomerName { get; set; } = string.Empty;
    public string Item { get; set; } = string.Empty;
    public decimal Amount { get; set; }
    public DateTime CreatedAt { get; set; } = DateTime.UtcNow;

    [JsonIgnore]
    public string Status { get; set; } = "Pending";
}

گام ۲: Producer تایپ‌شده با JSON Serialization

// Kafka/KafkaOrderProducer.cs
using System.Text.Json;
using Confluent.Kafka;
using KafkaTutorial.Api.Models;

namespace KafkaTutorial.Api.Kafka;

public class KafkaOrderProducer
{
    private readonly IProducer<Null, string> _producer;
    private readonly string _topic = KafkaConfig.TopicName;

    public KafkaOrderProducer()
    {
        var config = new ProducerConfig
        {
            BootstrapServers = KafkaConfig.BootstrapServers,
            Acks = Acks.All,  // منتظر تأیید همه Replicaها
            EnableIdempotence = true  // جلوگیری از ارسال دوباره
        };

        _producer = new ProducerBuilder<Null, string>(config).Build();
    }

    public async Task<(bool Success, string? Error)> SendOrderAsync(Order order)
    {
        try
        {
            var json = JsonSerializer.Serialize(order);
            var result = await _producer.ProduceAsync(
                _topic,
                new Message<Null, string> { Value = json }
            );

            Console.WriteLine(
                $"✅ سفارش ارسال شد — " +
                $"Order ID: {order.Id}, " +
                $"Partition: {result.Partition}, " +
                $"Offset: {result.Offset}"
            );

            return (true, null);
        }
        catch (Exception ex)
        {
            return (false, ex.Message);
        }
    }
}

توضیح:

  • Acks.All یعنی پیام حتماً روی چندتا سرور (Replica) ذخیره بشه قبل از تأیید.
  • EnableIdempotence = true یعنی اگه پیام به‌خاطر مشکل شبکه دوباره ارسال بشه، Kafka فقط یه نسخه رو نگه می‌داره.

گام ۳: Consumer با پردازش Order

// Kafka/KafkaOrderConsumer.cs
using System.Text.Json;
using Confluent.Kafka;
using KafkaTutorial.Api.Models;

namespace KafkaTutorial.Api.Kafka;

public class KafkaOrderConsumer : BackgroundService
{
    private readonly IServiceProvider _serviceProvider;
    private readonly ILogger<KafkaOrderConsumer> _logger;

    public KafkaOrderConsumer(
        IServiceProvider serviceProvider,
        ILogger<KafkaOrderConsumer> logger)
    {
        _serviceProvider = serviceProvider;
        _logger = logger;
    }

    protected override Task ExecuteAsync(CancellationToken stoppingToken)
    {
        return Task.Run(() => StartConsuming(stoppingToken), stoppingToken);
    }

    private void StartConsuming(CancellationToken cancellationToken)
    {
        var config = new ConsumerConfig
        {
            BootstrapServers = KafkaConfig.BootstrapServers,
            GroupId = "kitchen_service_group",
            AutoOffsetReset = AutoOffsetReset.Earliest,
            EnableAutoCommit = false  // خودمون Commit می‌کنیم
        };

        using var consumer = new ConsumerBuilder<Null, string>(config).Build();
        consumer.Subscribe(KafkaConfig.TopicName);

        _logger.LogInformation("🍽️ Kitchen Service شروع شد...");

        try
        {
            while (!cancellationToken.IsCancellationRequested)
            {
                var consumeResult = consumer.Consume(cancellationToken);

                try
                {
                    var order = JsonSerializer.Deserialize<Order>(consumeResult.Message.Value);

                    if (order != null)
                    {
                        _logger.LogInformation(
                            "🍳 آماده‌سازی سفارش: {OrderItem} برای {Customer} — مبلغ: {Amount}",
                            order.Item,
                            order.CustomerName,
                            order.Amount
                        );

                        // شبیه‌سازی پردازش
                        Thread.Sleep(500);

                        // Manual commit بعد از پردازش موفق
                        consumer.Commit(consumeResult);
                    }
                }
                catch (Exception ex)
                {
                    _logger.LogError(ex, "خطا در پردازش پیام: {Message}", consumeResult.Message.Value);
                    // اینجا می‌تونی پیام رو به یه Dead Letter Topic بفرستی
                }
            }
        }
        catch (OperationCanceledException)
        {
            _logger.LogInformation("Kitchen Service متوقف شد.");
        }
        finally
        {
            consumer.Close();
        }
    }
}

توضیح:

  • EnableAutoCommit = false یعنی خودمون بعد از پردازش موفق Offset رو Commit می‌کنیم. این خیلی مهمه!
  • اگه برنامه Crash کنه وسط پردازش، چون Commit نکردیم، پیام دوباره از همون Offset خونده می‌شه.
  • consumer.Commit(consumeResult) فقط بعد از موفقیت کامل اجرا می‌شه.

گام ۴: Notification Consumer (گروه متفاوت)

// Kafka/KafkaNotificationConsumer.cs
using System.Text.Json;
using Confluent.Kafka;
using KafkaTutorial.Api.Models;

namespace KafkaTutorial.Api.Kafka;

public class KafkaNotificationConsumer : BackgroundService
{
    private readonly ILogger<KafkaNotificationConsumer> _logger;

    public KafkaNotificationConsumer(ILogger<KafkaNotificationConsumer> logger)
    {
        _logger = logger;
    }

    protected override Task ExecuteAsync(CancellationToken stoppingToken)
    {
        return Task.Run(() => StartConsuming(stoppingToken), stoppingToken);
    }

    private void StartConsuming(CancellationToken cancellationToken)
    {
        var config = new ConsumerConfig
        {
            BootstrapServers = KafkaConfig.BootstrapServers,
            GroupId = "notification_service_group",  // گروه متفاوت!
            AutoOffsetReset = AutoOffsetReset.Earliest
        };

        using var consumer = new ConsumerBuilder<Null, string>(config).Build();
        consumer.Subscribe(KafkaConfig.TopicName);

        _logger.LogInformation("📧 Notification Service شروع شد...");

        try
        {
            while (!cancellationToken.IsCancellationRequested)
            {
                var consumeResult = consumer.Consume(cancellationToken);

                try
                {
                    var order = JsonSerializer.Deserialize<Order>(consumeResult.Message.Value);

                    _logger.LogInformation(
                        "📧 ایمیل تأیید ارسال شد به {Customer} برای سفارش: {Item}",
                        order?.CustomerName,
                        order?.Item
                    );
                }
                catch (Exception ex)
                {
                    _logger.LogError(ex, "خطا در Notification");
                }
            }
        }
        catch (OperationCanceledException)
        {
            _logger.LogInformation("Notification Service متوقف شد.");
        }
        finally
        {
            consumer.Close();
        }
    }
}

نکته کلیدی: هر دو Consumer (Kitchen و Notification) از یه Topic می‌خونن، چون GroupId متفاوت دارن، هرکدوم تمام پیام‌ها رو دریافت می‌کنن. این تفاوت اصلی با RabbitMQ است — چندتا Consumer مستقل می‌تونن همون پیام رو ببینن.

گام ۵: ثبت تو Program.cs

// Program.cs (نهایی)
using KafkaTutorial.Api.Kafka;

var builder = WebApplication.CreateBuilder(args);

builder.Services.AddSingleton<KafkaOrderProducer>();
builder.Services.AddHostedService<KafkaOrderConsumer>();
builder.Services.AddHostedService<KafkaNotificationConsumer>();

var app = builder.Build();

app.MapPost("/orders", async (KafkaOrderProducer producer, Order order) =>
{
    var (success, error) = await producer.SendOrderAsync(order);
    
    if (success)
        return Results.Ok(new { Message = "سفارش ثبت شد", OrderId = order.Id });
    
    return Results.StatusCode(500, new { Error = error });
});

app.Run();

گام ۶: تست

curl -X POST http://localhost:5000/orders \
  -H "Content-Type: application/json" \
  -d '{
    "customerName": "علی رضایی",
    "item": "پیتزا مارگاریتا",
    "amount": 120000
  }'

خروجی تو Console:

info: ✅ سفارش ارسال شد — Order ID: xxx, Partition: 0, Offset: 0
info: 🍽️ Kitchen Service شروع شد...
info: 📧 Notification Service شروع شد...
info: 🍳 آماده‌سازی سفارش: پیتزا مارگاریتا برای علی رضایی — مبلغ: 120000
info: 📧 ایمیل تأیید ارسال شد به علی رضایی برای سفارش: پیتزا مارگاریتا

۷. مباحث پیشرفته

۷.۱ Partitioning با Key

وقتی یه Key به پیام می‌دی، Kafka با Hash اون Key تصمیم می‌گیره پیام به کدوم Partition بره. این تضمین می‌کنه همه پیام‌های یه Key همیشه تو یه Partition باشن.

// Kafka/KafkaKeyedProducer.cs
using Confluent.Kafka;

namespace KafkaTutorial.Api.Kafka;

public class KafkaKeyedProducer
{
    private readonly IProducer<string, string> _producer;

    public KafkaKeyedProducer()
    {
        var config = new ProducerConfig
        {
            BootstrapServers = KafkaConfig.BootstrapServers
        };
        _producer = new ProducerBuilder<string, string>(config).Build();
    }

    public async Task SendWithKeyAsync(string key, string value)
    {
        var result = await _producer.ProduceAsync(
            KafkaConfig.TopicName,
            new Message<string, string> { Key = key, Value = value }
        );

        Console.WriteLine(
            $"Key: {key} → Partition: {result.Partition}, Offset: {result.Offset}"
        );
    }
}

توضیح: اگه ۳ تا Partition داشته باشی و Key "user-123" رو بفرستی، Kafka با Hash اون Key یه Partition مشخص (مثلاً Partition 1) انتخاب می‌کنه. همه پیام‌های "user-123" همیشه همون Partition 1 می‌رن. این برای ترتیب‌بندی پیام‌های یه کاربر حیاتی است.

۷.۲ Multiple Partitions و موازی‌سازی

برای ساخت Topic با چند Partition از CLI استفاده کن:

docker exec -it <container_id> kafka-topics.sh \
  --create \
  --bootstrap-server localhost:9092 \
  --topic orders \
  --partitions 3 \
  --replication-factor 1

توضیح: الان ۳ Partition داری. اگه ۳ Consumer تو یه Group داشته باشی، هرکدوم یه Partition رو می‌خونن و موازی کار می‌کنن.

۷.۳ Dead Letter Queue (DLQ)

وقتی پردازش یه پیام خطا می‌ده، به‌جای از دست دادنش، به یه Topic دیگه به اسم DLQ بفرست.

// Kafka/DlqHandler.cs
using Confluent.Kafka;

namespace KafkaTutorial.Api.Kafka;

public class DlqHandler
{
    private readonly IProducer<Null, string> _producer;
    private const string DlqTopic = "orders.DLQ";

    public DlqHandler()
    {
        var config = new ProducerConfig
        {
            BootstrapServers = KafkaConfig.BootstrapServers
        };
        _producer = new ProducerBuilder<Null, string>(config).Build();
    }

    public async Task SendToDlqAsync(string originalMessage, string error)
    {
        var dlqMessage = $"{originalMessage} | ERROR: {error}";
        await _producer.ProduceAsync(DlqTopic, new Message<Null, string> { Value = dlqMessage });
        Console.WriteLine($"⚠️ پیام به DLQ فرستاده شد: {error}");
    }
}

توضیح: DLQ یه الگوی استاندارد تو سیستم‌های Messaging است. بهت اجازه می‌ده پیام‌های مشکل‌دار رو بعداً بررسی و پردازش کنی، بدون اینکه کل Pipeline متوقف بشه.

۷.۴ Schema Registry و Avro

برای serialize پیام‌ها به‌صورت بهینه و version-safe، از Avro و Schema Registry استفاده می‌کنن.

dotnet add package Confluent.SchemaRegistry
dotnet add package Confluent.SchemaRegistry.Serdes.Avro
// مثال سریع (نیاز به Schema Registry Server داره)
var config = new ProducerConfig
{
    BootstrapServers = "localhost:9092"
};

var schemaRegistryConfig = new SchemaRegistryConfig
{
    Url = "http://localhost:8081"
};

var avroSerializerConfig = new AvroSerializerConfig
{
    AutoRegisterSchemas = true
};

توضیح: Avro یه فرمت باینری فشرده است. JSON خوانا ولی حجمش زیاده. Avro حجم کم و schema management قوی داره. برای پروژه‌های بزرگ، حتماً استفاده کن.

۷.۵ Exactly-Once Semantics

var config = new ProducerConfig
{
    BootstrapServers = KafkaConfig.BootstrapServers,
    EnableIdempotence = true,
    Acks = Acks.All,
    MaxInFlight = 5,
    MessageSendMaxRetries = int.MaxValue
};

توضیح: EnableIdempotence یعنی Kafka تضمین می‌کنه پیام دقیقاً یه‌بار ثبت بشه، حتی اگه شبکه دوباره ارسال کنه. این فقط برای Producer هست — برای Consumer باید خودت Idempotent بنویسی (مثلاً با چک کردن ID پیام).

۷.۶ Batch Processing

به‌جای ارسال یه‌به‌یه، می‌تونی batch‌بندی کنی:

public async Task SendBatchAsync(IEnumerable<string> messages)
{
    var tasks = messages.Select(msg => 
        _producer.ProduceAsync(KafkaConfig.TopicName, 
            new Message<Null, string> { Value = msg })
    );

    await Task.WhenAll(tasks);
}

توضیح: Batch کردن می‌تونه throughput رو ۱۰ برابر کنه. ولی مواظب باش — اگه batch خیلی بزرگ باشه، latency زیاد می‌شه.

۷.۷ Kafka Streams (ساده‌سازی)

Kafka Streams برای پردازش جریانی استفاده می‌شه. تو .NET می‌تونی از Streamiz.Kafka.Net استفاده کنی:

dotnet add package Streamiz.Kafka.Net
// مثال ساده Stream Processing
var config = new StreamConfig<StringSerDes, StringSerDes>
{
    ApplicationId = "order-stream-app",
    BootstrapServers = "localhost:9092"
};

var stream = new KafkaStream(config);
stream.Start();

توضیح: Kafka Streams یه API برای پردازش real-time داده‌هاست. می‌تونی aggregation، filtering، joins و... انجام بدی. به‌جای نوشتن Consumer دستی، از API سطح بالاتر استفاده می‌کنی.


۸. اشتباهات رایج

۸.۱ فعال کردن Auto Commit بی‌دقت

// ❌ اشتباه — Auto Commit فعاله
var config = new ConsumerConfig
{
    EnableAutoCommit = true,
    AutoCommitIntervalMs = 5000
};
// ✅ درست — Manual Commit
var config = new ConsumerConfig
{
    EnableAutoCommit = false
};
// بعد از پردازش موفق:
consumer.Commit(consumeResult);

چرا؟ با Auto Commit، ممکنه Offset قبل از پردازش کامل Commit بشه. اگه Crash کنی، پیام از دست می‌ره.

۸.۲ نداشتن تلاش مجدد (Retry)

وقتی پردازش خطا می‌ده، پیام رو رها نکن:

// ❌ اشتباه
catch (Exception ex)
{
    _logger.LogError(ex, "خطا");
    // ادامه می‌دی و پیام گم می‌شه
}
// ✅ درست
catch (Exception ex)
{
    _logger.LogError(ex, "خطا در پردازش، ارسال به DLQ");
    await _dlqHandler.SendToDlqAsync(consumeResult.Message.Value, ex.Message);
    consumer.Commit(consumeResult);  // بعد از ارسال به DLQ
}

۸.۳ تعداد Partition زیاد بدون دلیل

❌ 100 Partition برای 10 پیام در دقیقه
✅ 3 Partition کافیه

چرا؟ هر Partition یه file روی دیسک و یه overhead مدیریتی داره. Partition اضافی منابع رو هدر می‌ده.

۸.۴ Blocking در Consume

// ❌ اشتباه — پردازش طولانی تو حلقه
var result = consumer.Consume(cancellationToken);
Thread.Sleep(60000);  // پردازش طولانی
consumer.Commit(result);
// ✅ درست — پردازش async
var result = consumer.Consume(cancellationToken);
await ProcessAsync(result.Message.Value);
consumer.Commit(result);

۸.۵ یه Partition با چند Consumer تو یه Group

اگه ۱ Partition داشته باشی و ۳ Consumer تو یه Group، ۲ تا Consumer idle می‌شن. این اشتباهه!

۸.۶ Ignoring Partitioner

وقتی به ترتیب پیام‌های یه کاربر اهمیت داره، باید همیشه از Key استفاده کنی. بدون Key، ممکنه پیام‌های یه کاربر تو Partitionهای مختلف پخش بشن و ترتیب بهم بخوره.


۹. بهترین روش‌ها (Best Practices)

۹.۱ Use Dependency Injection

// ✅ اینجوری بنویس
builder.Services.AddSingleton<KafkaOrderProducer>();
builder.Services.AddHostedService<KafkaOrderConsumer>();

۹.۲ Use Configuration

// appsettings.json
{
  "Kafka": {
    "BootstrapServers": "localhost:9092",
    "TopicName": "orders",
    "ConsumerGroup": "order_processor_group"
  }
}
// Kafka/KafkaSettings.cs
public class KafkaSettings
{
    public string BootstrapServers { get; set; } = string.Empty;
    public string TopicName { get; set; } = string.Empty;
    public string ConsumerGroup { get; set; } = string.Empty;
}

// Program.cs
builder.Services.Configure<KafkaSettings>(builder.Configuration.GetSection("Kafka"));

۹.۳ همیشه Logging داشته باش

private readonly ILogger<KafkaOrderConsumer> _logger;

public async Task ProcessAsync(string message)
{
    _logger.LogInformation("شروع پردازش پیام: {Message}", message);
    try
    {
        // ...
        _logger.LogInformation("پردازش موفق");
    }
    catch (Exception ex)
    {
        _logger.LogError(ex, "خطا در پردازش");
        throw;
    }
}

۹.۴ Monitoring و Metrics

از Confluent Control Center، Kafka UI، یا Prometheus + Grafana استفاده کن.

# ساده‌ترین ابزار: Kafka UI
docker run -d -p 8080:8080 \
  -e KAFKA_CLUSTERS_0_NAME=local \
  -e KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS=host.docker.internal:9092 \
  provectuslabs/kafka-ui:latest

باز کن: http://localhost:8080

۹.۵ Partition Count رو از اول درست انتخاب کن

نرخ پیام Partition پیشنهادی
< 100/sec 1
100-1000/sec 3-6
1000-10000/sec 6-12
> 10000/sec 12+

نکته: بعد از ساخت Topic، اضافه کردن Partition امکان‌پذیره ولی ترتیب پیام‌ها بهم می‌خوره. پس از همون اول درست تصمیم بگیر.

۹.۶ Idempotent Consumer بنویس

public async Task ProcessOrderAsync(Order order)
{
    // چک کن آیا قبلاً پردازش شده؟
    if (await _orderRepository.ExistsAsync(order.Id))
    {
        _logger.LogInformation("سفارش قبلاً پردازش شده: {OrderId}", order.Id);
        return;
    }

    await _orderRepository.SaveAsync(order);
}

چرا؟ چون Kafka حداقل یه‌بار تحویل می‌ده. ممکنه یه پیام دوبار بیاد. اگه Consumerت Idempotent نباشه، یه سفارش دو بار ثبت می‌شه!

۹.۷ Graceful Shutdown

public class KafkaOrderConsumer : BackgroundService
{
    private IConsumer<Null, string>? _consumer;

    protected override Task ExecuteAsync(CancellationToken stoppingToken)
    {
        // ...
    }

    public override Task StopAsync(CancellationToken cancellationToken)
    {
        _consumer?.Close();  // Commit آخرین Offset و بستن تمیز
        return base.StopAsync(cancellationToken);
    }
}

۹.۸ Avro/Protobuf برای Production

JSON خوبه برای شروع، ولی تو Production از Avro یا Protobuf استفاده کن.


۱۰. منابع و ادامه مسیر

منابع رسمی

دوره‌ها و آموزش‌ها

ابزارهای مفید

  • Kafka UI — رابط گرافیکی برای مشاهده Topics، Consumers و...
  • Confluent Control Center — پلتفرم کامل مانیتورینگ
  • kafkacat (kcat) — CLI ابزار برای تست سریع
  • Schema Registry — مدیریت schemaهای Avro

مسیر ادامه

  1. اول: این آموزش رو با یه پروژه کوچک تمرین کن — یه Producer، یه Consumer، یه Topic.
  2. بعد: Schema Registry و Avro رو یاد بگیر.
  3. بعد: Kafka Streams رو بررسی کن.
  4. بعد: Kafka Connect برای اتصال به دیتابیس‌ها و سیستم‌های دیگه.
  5. بعد: ksqlDB برای query کردن جریان داده‌ها به‌صورت SQL.
  6. نهایتاً: Kafka Cluster Management و امنیت (SASL، SSL).

تمرین پیشنهادی

  1. یه پروژه ASP.NET Core با ۳ تا سرویس بساز:
    • OrderService (Producer)
    • KitchenService (Consumer Group A)
    • InventoryService (Consumer Group B)
  2. ۳ تا Partition بساز و موازی‌سازی رو تست کن.
  3. یه Consumer رو Stop کن و دوباره Start کن — باید از همون Offset ادامه بده.
  4. یه پیام خراب بفرست و DLQ رو امتحان کن.
  5. یه Dashboard ساده با یه API بزن که تعداد پیام‌های دریافتی رو نشون بده.

به پایان آموزش رسیدیم! اگه تا اینجا با دقت همراه بودی، الان می‌تونی یه سیستم Kafka واقعی با ASP.NET Core پیاده‌سازی کنی. مهم‌ترین چیز: شروع کن و تمرین کن. Kafka یه ابزار قدرتمنده، ولی پشتش یه منحنی یادگیری داره. قدم به قدم جلو بری، خیلی زود حرفه‌ای می‌شی.

Rejoining the server...

Rejoin failed... trying again in seconds.

Failed to rejoin.
Please retry or reload the page.

The session has been paused by the server.

Failed to resume the session.
Please retry or reload the page.